穷途末路的选择:Lambda 架构
0. 引言
有些计算(二度关联图谱、复杂统计学习/机器学习模型训练)因算法复杂度过高或数据量过大,无法直接实时计算。此时需要"压箱底"的 Lambda 架构:离线与实时相结合的折中方案。本文讲解 Lambda 架构的三层设计、其改进版 Kappa 架构,并用 Flink 演示实现。
1. Lambda 架构:三层各司其职
图表渲染中…
| 层 | 职责 | 存储与计算选型 |
|---|---|---|
| 批处理层 | 处理主数据集(历史全量数据) | HDFS/S3 + MapReduce/Hive/Spark;结果存 MySQL/HBase 供快速访问 |
| 快速处理层 | 处理增量数据(新进入系统的数据) | Flink/Spark Streaming/Storm;结果存 Redis 等高性能内存库 |
| 服务层 | 合并两层结果,提供查询服务 | 分别查询后合并返回 |
Lambda 架构是架构思想,各层技术选型不限定,可按实际情况组合。
主要问题:同一查询目标需为批处理层和快速处理层分别开发算法实现——相同逻辑两套代码、两个计算框架(如 Storm + Spark 并用),开发、测试、运维成本高。
2. Kappa 架构:用"流"统一编程界面
Kappa 架构的最大改进:用流计算技术取代批处理层——批处理层与快速处理层使用相同的流计算逻辑与统一的计算框架,降低开发测试运维成本。
存储也随之简化:不再将离线数据转储到 HDFS/S3 等块存储,而是按"流"的方式存于 Kafka 等流式存储系统,减少存储空间、避免数据转储、降低复杂度。
在 Flink、Spark Streaming 等流批一体计算框架与 Kafka、Pulsar 等流式存储系统的双重加持下,Kappa 已成为大数据处理的自然选择。
3. 用 Flink 实现 Kappa 架构
目标:统计"最近 3 天每种商品的销售量"。同一套代码,仅窗口类型不同:
java
DataStream counts = stream
.map(s -> JSONObject.parseObject(s, Event.class)) // 解析事件
.map(e -> new CountedEvent(e.product, 1, e.timestamp)) // 商品计数事件
.assignTimestampsAndWatermarks(new EventTimestampPeriodicWatermarks())
.keyBy("product")
// 批处理层:滑动窗口(3 天窗口、30 分钟步长)
.timeWindow(Time.days(3), Time.minutes(30))
// 快速处理层:翻滚窗口(15 秒)
// .window(TumblingEventTimeWindows.of(Time.seconds(15)))
.reduce((e1, e2) -> merge(e1, e2)); // 窗口内聚合批处理层与快速处理层的代码除窗口类型外完全相同——这正是 Kappa 架构的优势。
两层结果写入同一张数据库表(JdbcWriter 按层标记,ON DUPLICATE KEY UPDATE 幂等写入,并按步长取整对齐时间窗口便于合并);服务层一条 SQL 即可合并查询:
sql
SELECT product, sum(v_count) as s_count FROM (
SELECT * FROM table_counts WHERE start=? AND end=? AND layer='batch'
UNION
SELECT * FROM table_counts WHERE start>=? AND end<=? AND layer='fast'
) t GROUP BY product;4. 小结
- Lambda 架构:不能直接实时计算时,采用"离线(批处理层)+ 实时(快速处理层)+ 合并(服务层)"的折中方案;
- Kappa 架构:以"流"统一编程界面,同逻辑一套代码,配合流批一体框架与流式存储,显著降低复杂度;
- Kappa 不能完全取代 Lambda:批处理算法在丰富度、成熟度、第三方工具库方面仍优于流处理;是否统一框架,还需权衡数据人员(R/Python/Spark 等熟悉工具)与开发人员各自的产出效率;
- 最终方案应结合业务场景、技术积累、团队能力综合设计。
下一章讲解 Java 程序员的成长之路和从业方向。